跳到主要内容
Cowers://
全部文章
Python

Python 多线程与多进程:GIL、并发模型与实践

从 GIL 出发比较 threading、multiprocessing 与 concurrent.futures 的适用场景。

Python 并发模型总览:左侧单进程内的三个线程共享同一个 GIL,只有 Thread 1 在运行、另两个等待锁,三者共用一块 Heap;右侧两个进程各自持有独立 GIL 与独立内存、经 IPC 通信,因而真正并行。多线程适合 I/O 密集型(threading / asyncio),多进程适合 CPU 密集型(multiprocessing / ProcessPoolExecutor)。

ABSTRACT 概念速览

  • 多线程(threading):多个线程共享同一进程的内存空间;受 GIL 限制,适合 I/O 密集型 任务
  • 多进程(multiprocessing):每个进程有独立内存空间;绕过 GIL,适合 CPU 密集型 任务

一、GIL —— 理解并发的前提

NOTE 什么是 GIL? Global Interpreter Lock(全局解释器锁),CPython 的实现机制:同一时刻只有一个线程能执行 Python 字节码

  • I/O 操作(网络、磁盘)会释放 GIL → 多线程有效
  • 纯计算不释放 GIL → 多线程无法真正并行,用多进程代替
场景 推荐方案 原因
网络请求、文件读写 threading / asyncio I/O 等待期间释放 GIL
图像处理、数学计算 multiprocessing 独立进程,真正并行
两者都有 concurrent.futures 统一接口,灵活切换

二、threading —— 多线程

2.1 创建线程

import threading
 
def task(name, count):
    for i in range(count):
        print(f"[{name}] step {i}")
 
# 方式一:直接传函数
t = threading.Thread(target=task, args=("A", 3), kwargs={})
t.start()   # 启动线程
t.join()    # 等待线程结束(阻塞主线程)
 
# 方式二:继承 Thread 类(适合逻辑复杂时)
class MyThread(threading.Thread):
    def __init__(self, name):
        super().__init__()
        self.name = name
 
    def run(self):          # 重写 run(),不要直接调用 run()
        print(f"{self.name} running")
 
t = MyThread("Worker")
t.start()
t.join()

TIP start() vs run()

  • t.start() → 在新线程中执行 run()
  • t.run() → 在当前线程直接调用,没有新线程!

2.2 常用属性与方法

t = threading.Thread(target=task, args=("A",), daemon=True)
 
t.start()               # 启动
t.join(timeout=5.0)     # 最多等待 5 秒
t.is_alive()            # 是否还在运行 → bool
t.name                  # 线程名(可自定义)
t.daemon                # 守护线程:主线程结束时强制退出
 
threading.current_thread()   # 返回当前线程对象
threading.enumerate()        # 返回所有存活线程列表
threading.active_count()     # 存活线程数量

WARNING 守护线程(daemon) daemon=True 时,主线程结束,守护线程立即被杀死,不等其完成。
适合后台心跳、日志刷新等"可抛弃"任务。

2.3 线程锁(Lock)

多线程共享内存 → 需要加锁防止竞态条件

import threading
 
counter = 0
lock = threading.Lock()
 
def increment():
    global counter
    with lock:          # 推荐:自动 acquire / release
        counter += 1
 
# 等价于:
def increment_manual():
    global counter
    lock.acquire()
    try:
        counter += 1
    finally:
        lock.release()  # 必须放在 finally!
 
threads = [threading.Thread(target=increment) for _ in range(1000)]
for t in threads: t.start()
for t in threads: t.join()
print(counter)  # 1000(无锁时结果不确定)

2.4 RLock —— 可重入锁

rlock = threading.RLock()
 
def outer():
    with rlock:
        inner()     # 同一线程可再次获取,不会死锁
 
def inner():
    with rlock:     # Lock 在这里会死锁,RLock 不会
        print("inner")

2.5 Semaphore —— 信号量(限制并发数)

sem = threading.Semaphore(3)   # 最多 3 个线程同时运行
 
def limited_task(i):
    with sem:
        print(f"Task {i} running")
        time.sleep(1)
 
threads = [threading.Thread(target=limited_task, args=(i,)) for i in range(10)]
for t in threads: t.start()
for t in threads: t.join()

2.6 Event —— 线程间通信

event = threading.Event()
 
def waiter():
    print("等待信号...")
    event.wait()            # 阻塞,直到 event 被 set()
    print("收到信号,继续!")
 
def setter():
    time.sleep(2)
    event.set()             # 通知所有等待的线程
 
threading.Thread(target=waiter).start()
threading.Thread(target=setter).start()
方法 说明
event.set() 设置标志,释放所有 wait()
event.clear() 清除标志,后续 wait() 会再次阻塞
event.is_set() 返回当前状态
event.wait(timeout) 阻塞等待,可设超时

2.7 Queue —— 线程安全队列

from queue import Queue
 
q = Queue(maxsize=5)    # 0 表示无限制
 
q.put(item)             # 放入,队满则阻塞
q.put_nowait(item)      # 队满直接抛 Full 异常
q.get()                 # 取出,队空则阻塞
q.get_nowait()          # 队空直接抛 Empty 异常
q.task_done()           # 告知该任务已完成
q.join()                # 阻塞,直到所有 task_done() 被调用
q.qsize()               # 当前大小
q.empty() / q.full()    # 状态判断(非原子,仅供参考)

生产者-消费者模式:

from queue import Queue
import threading, time
 
q = Queue()
 
def producer():
    for i in range(5):
        q.put(i)
        print(f"生产: {i}")
        time.sleep(0.5)
 
def consumer():
    while True:
        item = q.get()
        if item is None:    # 哨兵值:退出信号
            break
        print(f"消费: {item}")
        q.task_done()
 
t1 = threading.Thread(target=producer)
t2 = threading.Thread(target=consumer)
t1.start(); t2.start()
t1.join()
q.put(None)     # 发送退出信号
t2.join()

三、multiprocessing —— 多进程

3.1 创建进程

from multiprocessing import Process
 
def task(name):
    print(f"进程 {name} 运行中,PID={os.getpid()}")
 
if __name__ == "__main__":     # ⚠️ Windows 必须有这个守卫!
    p = Process(target=task, args=("Worker",))
    p.start()
    p.join()

DANGER if __name__ == "__main__" 守卫 Windows 使用 spawn 方式创建进程,会重新 import 主模块。
没有这个守卫会无限递归创建进程! Linux/macOS 用 fork,风险较小但仍建议加上。

3.2 常用属性与方法

p = Process(target=task, args=("A",), daemon=False)
 
p.start()
p.join(timeout=5)
p.terminate()       # 发送 SIGTERM(较优雅)
p.kill()            # 发送 SIGKILL(立即杀死)
p.is_alive()
p.pid               # 进程 ID
p.exitcode          # 退出码(None 表示还在运行)

3.3 进程间通信(IPC)

进程内存相互隔离,需要专门的 IPC 机制:

Queue(多进程安全)

from multiprocessing import Process, Queue
 
def producer(q):
    q.put("hello from producer")
 
def consumer(q):
    msg = q.get()
    print(f"收到: {msg}")
 
if __name__ == "__main__":
    q = Queue()
    p1 = Process(target=producer, args=(q,))
    p2 = Process(target=consumer, args=(q,))
    p1.start(); p2.start()
    p1.join(); p2.join()

Pipe(双向管道)

from multiprocessing import Process, Pipe
 
def child(conn):
    conn.send("来自子进程")
    msg = conn.recv()
    print(f"子进程收到: {msg}")
    conn.close()
 
if __name__ == "__main__":
    parent_conn, child_conn = Pipe()    # 返回两端
    p = Process(target=child, args=(child_conn,))
    p.start()
    print(f"父进程收到: {parent_conn.recv()}")
    parent_conn.send("来自父进程")
    p.join()

共享内存(Value / Array)

from multiprocessing import Process, Value, Array
import ctypes
 
def modify(val, arr):
    val.value = 42
    arr[0] = 100
 
if __name__ == "__main__":
    v = Value(ctypes.c_int, 0)         # 共享整数
    a = Array(ctypes.c_int, [1,2,3])   # 共享数组
 
    p = Process(target=modify, args=(v, a))
    p.start()
    p.join()
    print(v.value, list(a))   # 42 [100, 2, 3]

NOTE 共享内存 vs Queue/Pipe

  • Value/Array:零拷贝,高性能,适合大量数值数据
  • Queue/Pipe:序列化传输,更灵活,适合任意 Python 对象

3.4 进程锁

from multiprocessing import Process, Lock
 
def task(lock, i):
    with lock:
        print(f"进程 {i} 安全写入")
 
if __name__ == "__main__":
    lock = Lock()
    processes = [Process(target=task, args=(lock, i)) for i in range(5)]
    for p in processes: p.start()
    for p in processes: p.join()

3.5 进程池(Pool)

TIP 何时用 Pool? 需要执行大量同类任务时,Pool 自动管理进程复用,避免频繁创建销毁进程的开销。

from multiprocessing import Pool
 
def square(x):
    return x * x
 
if __name__ == "__main__":
    with Pool(processes=4) as pool:   # 4 个工作进程
 
        # map:阻塞,返回有序结果列表
        result = pool.map(square, range(10))
        print(result)   # [0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
 
        # apply:同步执行单个任务
        r = pool.apply(square, args=(5,))
        print(r)   # 25
 
        # apply_async:异步执行单个任务
        async_result = pool.apply_async(square, args=(7,))
        print(async_result.get(timeout=5))   # 49
 
        # map_async:异步 map
        async_results = pool.map_async(square, range(5))
        print(async_results.get())   # [0, 1, 4, 9, 16]
 
        # imap:惰性迭代(大数据节省内存)
        for r in pool.imap(square, range(5)):
            print(r)
方法 是否阻塞 返回值 适用场景
map(f, iterable) ✅ 阻塞 list 同步批量处理
map_async(f, iterable) ❌ 非阻塞 AsyncResult 异步批量处理
apply(f, args) ✅ 阻塞 单个结果 同步单次调用
apply_async(f, args) ❌ 非阻塞 AsyncResult 异步单次调用
imap(f, iterable) ✅ 惰性 迭代器 大数据流式处理
starmap(f, iterable) ✅ 阻塞 list 多参数展开
# starmap:多参数
def add(a, b): return a + b
 
with Pool(4) as pool:
    result = pool.starmap(add, [(1,2), (3,4), (5,6)])
    print(result)   # [3, 7, 11]

四、concurrent.futures —— 统一高层接口

ABSTRACT 核心优势同一套 API 切换多线程和多进程,代码改动极小。

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
 
def task(n):
    return n * n
 
# 多线程(I/O 密集型)
with ThreadPoolExecutor(max_workers=4) as executor:
    futures = [executor.submit(task, i) for i in range(5)]
    for f in futures:
        print(f.result())   # 阻塞等待结果
 
# 多进程(CPU 密集型)—— 只改一行!
with ProcessPoolExecutor(max_workers=4) as executor:
    futures = [executor.submit(task, i) for i in range(5)]
    for f in futures:
        print(f.result())

map 方式

with ThreadPoolExecutor(max_workers=4) as executor:
    # 按提交顺序返回结果
    results = list(executor.map(task, range(10)))
    print(results)

Future 对象

future = executor.submit(task, 42)
 
future.result(timeout=5)    # 获取结果,超时抛 TimeoutError
future.done()               # 是否完成
future.cancelled()          # 是否被取消
future.cancel()             # 尝试取消(若已开始执行则无效)
future.exception()          # 若任务异常,返回异常对象

as_completed —— 哪个先完成先处理

from concurrent.futures import as_completed
import time, random
 
def slow_task(n):
    time.sleep(random.random())
    return n * 2
 
with ThreadPoolExecutor(max_workers=4) as executor:
    futures = {executor.submit(slow_task, i): i for i in range(5)}
    for future in as_completed(futures):    # 先完成先返回
        n = futures[future]
        print(f"任务 {n} 完成,结果: {future.result()}")

五、对比速查表

多线程 vs 多进程

维度 threading multiprocessing
内存 共享 独立
GIL 影响 受限(I/O 可并发) 不受限(真正并行)
适合场景 I/O 密集 CPU 密集
通信方式 共享变量 + 锁 Queue / Pipe / 共享内存
创建开销
稳定性 崩溃影响整个进程 进程间隔离

同步原语速查

原语 模块 作用
Lock threading / multiprocessing 互斥锁,同一时刻只有一个持有者
RLock threading 可重入锁,同一线程可多次获取
Semaphore threading / multiprocessing 限制并发数量
Event threading / multiprocessing 线程/进程间信号通知
Condition threading 复杂条件同步(wait/notify)
Barrier threading 等待所有线程到达同一屏障点

六、常见陷阱

DANGER 陷阱 1:多进程忘写 __main__ 守卫

# ❌ Windows 下无限递归
p = Process(target=f)
p.start()
 
# ✅ 正确
if __name__ == "__main__":
    p = Process(target=f)
    p.start()

DANGER 陷阱 2:线程中修改共享变量不加锁

# ❌ 竞态条件,结果不确定
counter += 1
 
# ✅ 加锁保护
with lock:
    counter += 1

DANGER 陷阱 3:多进程中普通全局变量不共享

count = 0   # 每个进程有自己的副本,修改互不影响
 
# ✅ 需要共享时用 Value / Manager
from multiprocessing import Value
count = Value('i', 0)

DANGER 陷阱 4:子线程异常不会传播到主线程

# ❌ 子线程抛异常,主线程不知道
t = Thread(target=risky_func)
t.start(); t.join()   # 不会抛异常!
 
# ✅ 用 concurrent.futures 会把异常保存到 Future
future = executor.submit(risky_func)
future.result()   # 这里才会重新抛出异常

WARNING 陷阱 5:死锁 多个锁按不同顺序获取时可能死锁:

# 线程 1: 先拿 lock_a 再拿 lock_b
# 线程 2: 先拿 lock_b 再拿 lock_a
# → 互相等待,永远阻塞
 
# ✅ 统一加锁顺序,或使用 timeout
lock_a.acquire(timeout=1)

七、选型决策树

需要并发?
├─ I/O 密集(网络/文件)
│   ├─ 任务数不多 → threading.Thread
│   ├─ 任务数很多 → ThreadPoolExecutor
│   └─ 大量异步 I/O → asyncio(参见异步笔记)

└─ CPU 密集(计算)
    ├─ 任务数不多 → multiprocessing.Process
    ├─ 任务数很多 → ProcessPoolExecutor 或 Pool
    └─ 需要共享大量数据 → multiprocessing + shared memory

关联笔记:Python异步编程 · Python GIL深入 · asyncio事件循环